SPB Git

spb/search-box Public

Agentic web research engine — hypotheses, verbatim evidence, contradictions, sourced answers streamed live. Claude Opus 5 + Firecrawl + PostgreSQL.

TypeScript 76.9% CSS 18.7% SQL 2.1% JavaScript 1.8% Shell 0.5%
2.8 KB · 96 lines typescript
Raw Blame History
1/**2 * Search-box.ai3 * Author: Simon-Pierre Boucher4 * Contact: contact@spboucher.ai5 * File: apps/web/app/api/research/[id]/stream/route.ts6 * Description: SSE stream of research events with Last-Event-ID replay (proxy-safe headers).7 */89import { events, sessions } from "@search-box/db";10import { ensureMigrated } from "@/lib/runner";1112export const runtime = "nodejs";13export const dynamic = "force-dynamic";1415const POLL_MS = 400;16const HEARTBEAT_MS = 15_000;1718export async function GET(19  req: Request,20  { params }: { params: Promise<{ id: string }> }21): Promise<Response> {22  await ensureMigrated();23  const { id } = await params;24  const session = await sessions.get(id);25  if (!session) return new Response("not found", { status: 404 });2627  // Replay position: Last-Event-ID header (EventSource reconnect) or ?from=28  const url = new URL(req.url);29  const fromParam = url.searchParams.get("from");30  const lastEventId = req.headers.get("last-event-id");31  let cursor = Number(lastEventId ?? fromParam ?? 0);32  if (!Number.isFinite(cursor) || cursor < 0) cursor = 0;3334  const encoder = new TextEncoder();3536  const stream = new ReadableStream<Uint8Array>({37    async start(controller) {38      let closed = false;39      const close = () => {40        if (!closed) {41          closed = true;42          try {43            controller.close();44          } catch {45            /* already closed */46          }47        }48      };49      req.signal.addEventListener("abort", close);5051      const send = (seq: number, data: string) => {52        controller.enqueue(encoder.encode(`id: ${seq}\ndata: ${data}\n\n`));53      };5455      let lastBeat = Date.now();56      try {57        while (!closed) {58          const batch = await events.listAfter(id, cursor);59          for (const ev of batch) {60            cursor = ev.seq;61            send(ev.seq, JSON.stringify(ev.payload));62          }6364          // Stop once the session reached a terminal state and all events were sent.65          if (batch.length === 0) {66            const current = await sessions.get(id);67            if (current && (current.status === "completed" || current.status === "failed")) {68              controller.enqueue(encoder.encode(`event: end\ndata: {}\n\n`));69              break;70            }71          }7273          if (Date.now() - lastBeat >= HEARTBEAT_MS) {74            controller.enqueue(encoder.encode(`: heartbeat\n\n`));75            lastBeat = Date.now();76          }77          await new Promise((r) => setTimeout(r, POLL_MS));78        }79      } catch {80        // client disconnected or db error — just end the stream81      } finally {82        close();83      }84    }85  });8687  return new Response(stream, {88    headers: {89      "Content-Type": "text/event-stream; charset=utf-8",90      "Cache-Control": "no-cache, no-transform",91      Connection: "keep-alive",92      "X-Accel-Buffering": "no"93    }94  });95}96